package learn

import domain.Domain.Access
import org.apache.flink.api.scala._
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment

object test01Stage {
    def main(args: Array[String]): Unit = {
        val env = StreamExecutionEnvironment.getExecutionEnvironment
        env.setParallelism(1)
        env.socketTextStream("hadoop000", 9527).map(x=>{
            val split = x.split(",")
            //Access(split(0), split(1))
            (split(0), 1)
        }).keyBy(0).sum(1).setParallelism(4).print()

        env.execute("stageTest")
    }
}
